Skip to content

Kafka Connect: Rework commit-coordinator leader election and harden the coordinator - #17450

Open
kumarpritam863 wants to merge 2 commits into
apache:mainfrom
kumarpritam863:feature/hardening-plus-election
Open

Kafka Connect: Rework commit-coordinator leader election and harden the coordinator#17450
kumarpritam863 wants to merge 2 commits into
apache:mainfrom
kumarpritam863:feature/hardening-plus-election

Conversation

@kumarpritam863

Copy link
Copy Markdown
Contributor

What

Reworks how the Iceberg Kafka Connect sink elects its single commit Coordinator and
hardens the coordinator's fencing, recovery, and shutdown paths. No control-topic
wire-format change; exactly-once semantics are preserved.

Leader election — level read, no Admin call

  • Before: the coordinator was chosen on every rebalance by calling
    Admin.describeConsumerGroups(connectGroupId) and picking the task that owned the
    globally-lowest (topic, partition) from a transient, mid-rebalance member snapshot.
    This required a DESCRIBE ACL, added an Admin round-trip to the rebalance path, and was
    racy during cooperative rebalancing.
  • After: the leader is the task whose assignment() contains partition 0 of the
    lexicographically-smallest topic in its own subscription()
    — connector-wide,
    identical across tasks, and valid for both a topics list and a topics.regex. No Admin
    call.
  • Leadership is read as a level on the task thread. open/close only flag that a
    reconcile is needed; save() reconciles once — start the coordinator if this task owns
    the leader partition, stop it otherwise (both idempotent). Keeping start/stop off the
    rebalance callback avoids blocking RPCs there and stops an eager rebalance from
    needlessly restarting a still-leading coordinator (which would discard in-flight commit
    state).
  • The coordinator's commit-readiness partition count now comes from
    consumer.partitionsFor(...) over the subscribed topics instead of a member snapshot.
  • Adds a Committer.configure(Catalog, IcebergSinkConfig, SinkTaskContext) lifecycle hook
    (invoked from IcebergSinkTask.start) for one-time setup.

Coordinator hardening

  • Zombie fencing via a fixed transactional.id. The coordinator producer id is now
    <connectGroupId>-<connectorName>-coord — identical across a connector's tasks and
    stable across restarts — so a newly elected coordinator's initTransactions()
    epoch-fences a prior (zombie) coordinator's control-plane writes. Worker ids are
    unchanged. This fences the brief two-coordinator overlap the level-read handoff can
    create (the losing task stops on its next save()).
  • Fenced ≠ fatal. A coordinator that terminates because it was fenced
    (ProducerFenced / InvalidProducerEpoch / UnknownProducerId, matched across the
    cause chain) is cleared without failing the task; any other termination still fails the
    task.
  • Recovery reads from earliest. The -coord consumer group defaults to
    auto.offset.reset=earliest, so a fresh or expired group re-reads uncommitted control
    events instead of skipping to the log end. Replay is idempotent via the snapshot offset
    floor + distinctByKey(location) dedup + the offsets compare-and-swap.
  • Idempotency filter in Channel.consumeAvailable. Control-topic records at or below
    the already-consumed per-partition offset are skipped, so a re-delivered or rewound
    record is never re-buffered or re-counted.
  • Transient Kafka commit errors are retryable. Errors from commitConsumerOffsets
    (Iceberg CommitFailedException, Kafka CommitFailedException,
    RebalanceInProgressException, RetriableException) are retried rather than failing the
    task — the table commit already succeeded.
  • Bounded, interrupt-safe shutdown. stopCoordinator clears state first, then
    bounded-joins the thread; terminate() failures are best-effort and interrupts are
    preserved.
  • No producer leak on failed init. KafkaClientFactory.createProducer closes the
    producer if initTransactions() throws.

Compatibility

  • No config or control-topic wire-format changes.
  • Exactly-once unchanged — anchored on the Iceberg offsets compare-and-swap +
    file-location dedup; the fixed transactional.id adds control-plane fencing on top.

Testing

  • Unit tests for the election key (leaderPartition), the retryable-commit classification,
    and coordinator fenced-vs-fatal termination.
  • Integration suite passes (13 tests).
  • Raised the integration-test commit-wait from 30s to 60s to reduce CI-load flakiness.

Pritam Kumar Mishra added 2 commits July 31, 2026 17:09
…he coordinator

Reworks how the sink elects its single commit Coordinator and hardens the
coordinator's fencing, recovery, and shutdown paths. No control-topic wire-format
change; exactly-once semantics are preserved.

Leader election (level read, no Admin call)
- Before: the coordinator was chosen on every rebalance by calling
  Admin.describeConsumerGroups(connectGroupId) and picking the task that owned the
  globally-lowest (topic, partition) from a transient, mid-rebalance member snapshot.
  This needed a DESCRIBE ACL, added an Admin round-trip to the rebalance path, and
  was racy during cooperative rebalancing.
- After: the leader is the task whose assignment contains partition 0 of the
  lexicographically-smallest topic in its own consumer subscription() -- connector-wide,
  identical across tasks, and valid for both a topics list and a topics.regex. No Admin
  call. Leadership is read as a level on the task thread: open()/close() only flag that a
  reconcile is needed, and save() reconciles once (start the coordinator if this task
  owns the leader partition, stop it otherwise). Keeping start/stop off the rebalance
  callback avoids blocking RPCs there and stops an eager rebalance from needlessly
  restarting a still-leading coordinator (which would discard in-flight commit state).
- The coordinator's commit-readiness partition count is derived from
  consumer.partitionsFor(...) over the subscribed topics instead of a member-assignment
  snapshot.
- Adds a Committer.configure(Catalog, IcebergSinkConfig, SinkTaskContext) lifecycle hook
  (invoked from IcebergSinkTask.start) for one-time setup.

Coordinator hardening
- Zombie fencing via a fixed transactional.id. The coordinator producer id is now
  connectGroupId-connectorName-coord -- identical across a connector's tasks and stable
  across restarts -- so a newly elected coordinator's initTransactions() epoch-fences a
  prior (zombie) coordinator's control-plane writes. Worker ids are unchanged. This
  fences the brief two-coordinator overlap that the level-read handoff can create (the
  losing task stops on its next save()).
- Fenced != fatal. A coordinator that terminates because it was fenced
  (ProducerFenced / InvalidProducerEpoch / UnknownProducerId, matched across the cause
  chain) is cleared without failing the task; any other termination still fails the task.
- Recovery reads from earliest. The coordinator's -coord consumer group defaults to
  auto.offset.reset=earliest, so a fresh or expired group re-reads uncommitted control
  events instead of skipping to the log end. Replay is idempotent via the snapshot offset
  floor + distinctByKey(location) dedup + the offsets compare-and-swap.
- Idempotency filter in Channel.consumeAvailable: control-topic records at or below the
  already-consumed per-partition offset are skipped, so a re-delivered or rewound record
  is never re-buffered or re-counted.
- Transient Kafka commit errors from commitConsumerOffsets (Iceberg CommitFailedException,
  Kafka CommitFailedException, RebalanceInProgressException, RetriableException) are
  retried rather than failing the task; the table commit already succeeded.
- Bounded, interrupt-safe shutdown: stopCoordinator clears state first, then bounded-joins
  the thread; terminate() failures are best-effort and interrupts are preserved.
- KafkaClientFactory.createProducer closes the producer if initTransactions() throws.

Compatibility
- No config or control-topic wire-format changes.
- Exactly-once unchanged -- anchored on the Iceberg offsets compare-and-swap +
  file-location dedup; the fixed transactional.id adds control-plane fencing on top.
- Known edge (OCC-safe): with topics.regex, a lexicographically-smaller topic appearing
  later can briefly run two coordinators if the losing task receives no rebalance callback.

Testing
- Unit tests for the election key (leaderPartition), the retryable-commit classification,
  and coordinator fenced-vs-fatal termination; integration suite (13 tests) passes.
- Raised the integration-test commit-wait from 30s to 60s to reduce CI-load flakiness.
@kumarpritam863

Copy link
Copy Markdown
Contributor Author

@laskoviymishka @danicafine @bryanck can you please review.

@laskoviymishka
laskoviymishka self-requested a review August 1, 2026 10:02
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant